在需要很多執行緒幫忙處理工作時,如果每次都開一個新執行緒、用完再銷毀顯得非常浪費,不如直接準備好一批執行緒,看誰有空就請誰做事,這批「先被準備好、待命中的執行緒」就稱為「執行緒池(thread pool)」。
以下 C++ 程式碼展示了執行緒池的運作方式:
#include <iostream>
#include <thread>
#include <vector>
#include <queue>
#include <mutex>
#include <condition_variable>
#include <functional>
#include <chrono>
class ThreadPool {
private:
std::queue<std::function<void()>> tasks;
std::vector<std::thread> workers;
std::mutex mtx;
std::mutex cout_mtx;
std::condition_variable cv;
bool done = false;
void worker(int id) {
while (true) {
std::function<void()> task;
{
std::unique_lock<std::mutex> lock(mtx);
cv.wait(lock, [this] { return !tasks.empty() || done; });
if (done && tasks.empty()) {
return;
}
task = tasks.front();
tasks.pop();
}
{
std::lock_guard<std::mutex> lock(cout_mtx);
std::cout << "Worker " << id << " starts a task\n";
}
task();
}
}
public:
ThreadPool(int numWorkers) {
for (int i = 0; i < numWorkers; i++) {
workers.emplace_back(&ThreadPool::worker, this, i);
}
}
void addTask(std::function<void()> task) {
{
std::lock_guard<std::mutex> lock(mtx);
tasks.push(task);
}
cv.notify_one();
}
void stop() {
{
std::lock_guard<std::mutex> lock(mtx);
done = true;
}
cv.notify_all();
}
void wait() {
for (auto& t : workers) {
t.join();
}
}
std::mutex& getCoutMutex() {
return cout_mtx;
}
~ThreadPool() {
stop();
wait();
}
};
int main() {
ThreadPool pool(3);
for (int i = 1; i <= 6; i++) {
pool.addTask([i, &pool] {
{
std::lock_guard<std::mutex> lock(pool.getCoutMutex());
std::cout << " Task " << i << " is running\n";
}
std::this_thread::sleep_for(std::chrono::seconds(1));
{
std::lock_guard<std::mutex> lock(pool.getCoutMutex());
std::cout << " Task " << i << " is done\n";
}
});
}
return 0;
}
以下透過執行順序,逐步說明上述程式碼內容。
透過 ThreadPool 類別建立實例,名為 pool、傳入參數 3 代表預計要在執行緒池建立三個執行緒。
// main() 函數中
ThreadPool pool(3);
建立 pool 實例會觸發 ThreadPool 類別中位於 public 的建構子 ThreadPool:
// ThreadPool 類別 public 區建構子
ThreadPool(int numWorkers) { // 傳入的參數為 3
for (int i = 0; i < numWorkers; i++) {
// 下方將詳細說明下面這行
workers.emplace_back(&ThreadPool::worker, this, i);
}
}
workers 來自定義類別時寫在 private 區塊的工人陣列,裡面每個元素都要裝執行緒(每個執行緒就是工人):
// ThreadPool 類別 private 區
std::vector<std::thread> workers;
這個陣列要用 emplace_back() 加入一個一個的工人元素,之所以用 emplace_back() 而不用 push_back(),是為了讓每個元素都直接在 workers 陣列裡面就地建造,而非建立臨時物件再複製或搬移到 workers 陣列中。
後面的 &ThreadPool::worker, this 表示要執行的是 ThreadPool 本身物件中的 worker() 函數,而 i 就是迴圈要跑的 0、1、2,因此整體:
// ThreadPool 類別 public 區建構子
workers.emplace_back(&ThreadPool::worker, this, i);
可以想成:
workers 陣列中,this -> worker(i)
因為 numWorkers 參數為 3,所以:
i = 0:開一條執行緒跑 this -> worker(0)
i = 1:開一條執行緒跑 this -> worker(1)
i = 2:開一條執行緒跑 this -> worker(2)
而這個 worker() 函數做了很多事。
首先是三個 id 分別為 0、1、2 的執行緒都在無窮迴圈中不斷尋找工作,如果有拿到工作,就放到 task 變數中:
// ThreadPool 類別 private 區
void worker(int id) {
while (true) {
每一項工作都會是沒有回傳型別、也不會有參數的函數 function<void()>:
// worker() 函數中
std::function<void()> task;
每個執行緒都會試著搶互斥鎖:
// worker() 函數中
std::unique_lock<std::mutex> lock(mtx);
但此時還沒有工作要做、tasks 為空,且 done 為 false,因此三個執行緒都會釋放鎖、進入沉睡狀態,等待條件變數 cv 的訊號:
// worker() 函數中
cv.wait(lock, [this] { return !tasks.empty() || done; });
截至目前為止,三個執行緒(工人)都已經準備好要工作,就等實際工作下來。
在 private 區有先建立一個用來裝工作的工作佇列 tasks:
// ThreadPool 類別 private 區
std::queue<std::function<void()>> tasks;
在 public 區也有一個將 tasks 工作佇列增添工作的函數 addTask():
// ThreadPool 類別 public 區
void addTask(std::function<void()> task) { // 接 task 函數為參數
{
std::lock_guard<std::mutex> lock(mtx); // 把工作佇列 tasks 上鎖保護起來
tasks.push(task); // 把 task 加進 tasks 工作佇列
}
cv.notify_one(); // 喚醒一個沉睡中的執行緒
}
在 main() 函數中透過 for 迴圈生產六個任務:
// main() 函數中
for (int i = 1; i <= 6; i++)
而針對每個任務,都執行 pool 物件實例中的 addTask() 方法:
// main() 函數中
pool.addTask(
// ...
);
剛剛已經看到 addTask() 函數接的參數是 std::function<void()> task,而這個 task 函數長這樣:
[i, &pool] { // i 表示迴圈跑到幾(建立第幾個工作)、&pool 表示可以取用 pool 實例物件中定義的 getCoutMutex() 方法
{
std::lock_guard<std::mutex> lock(pool.getCoutMutex()); // 先為 cout 上鎖
std::cout << " Task " << i << " is running\n"; // 模擬開始工作
}
std::this_thread::sleep_for(std::chrono::seconds(1)); // 模擬工作了一秒
{
std::lock_guard<std::mutex> lock(pool.getCoutMutex()); // 為 cout 上鎖
std::cout << " Task " << i << " is done\n"; // 模擬完成工作
}
} // 離開大括號時自動釋放鎖
addTask() 函數會把第 i 個 task 函數塞到 tasks 工作佇列中:
// addTask() 函數中
tasks.push(task);
如果 tasks 被塞進第 1 個 task,那 tasks 佇列就會變成:
[ [1, &pool]{std:: ... is done\n";} ]
隨著迴圈由 1 跑到 6,tasks 這個工作佇列會逐漸變成:
[ [1, &pool]{std:: ... is done\n";}, [2, &pool]{std:: ... is done\n";}, ... , [6, &pool]{std:: ... is done\n";}]
↑ 第 0 個元素,i 為 1 ↑ 第 1 個元素,i 為 2 ↑ 第 5 個元素,i 為 6
addTask() 函數離開大括號為工作佇列 tasks 解鎖後,還會去喚醒一個沉睡中的執行緒:
// addTask() 函數中
cv.notify_one();
其中一個工人(執行緒)被喚醒,去搶 mtx 鎖:
// worker() 函數中
std::unique_lock<std::mutex> lock(mtx);
此時 tasks 佇列已經不為空,因此 !tasks.empty() 為 true,不會原地休息而會往下執行:
// worker() 函數中
cv.wait(lock, [this] { return !tasks.empty() || done; }); // 不會卡在這裡休息
因為 tasks.empty() 不為 true,因此跳過這段(不會直接 return):
// worker() 函數中
if (done && tasks.empty()) {
return;
}
每個 worker 函數都透過以下程式碼,被喚醒後從最前面的工作開始取:
// worker() 函數中
task = tasks.front();
tasks.pop();
假設原本的工作佇列 tasks 長這樣:
[ [1, &pool]{std:: ... is done\n";}, [2, &pool]{std:: ... is done\n";}, ... , [6, &pool]{std:: ... is done\n";}]
那經過上兩行程式碼後:
task 為 [1, &pool]{std:: ... is done\n";}。tasks 變成:[ [2, &pool]{std:: ... is done\n";}, [3, &pool]{std:: ... is done\n";}, ... , [6, &pool]{std:: ... is done\n";}]
接著顯示被喚醒的工人要開始工作了:
// worker() 函數中
{
std::lock_guard<std::mutex> lock(cout_mtx);
std::cout << "Worker " << id << " starts a task\n";
}
透過呼叫 task() 函數,真正執行工作內容(也就是寫在 main() 函數中 lambda 函數那些):
// worker() 函數中
task();
假設是第 1 個工作被執行,那就是跑這個:
// main() 函數中
[1, &pool] {
{
std::lock_guard<std::mutex> lock(pool.getCoutMutex()); // 先為 cout 上鎖
std::cout << " Task " << i << " is running\n"; // 這裡的 i 以 1 代替
}
std::this_thread::sleep_for(std::chrono::seconds(1));
{
std::lock_guard<std::mutex> lock(pool.getCoutMutex()); // 為 cout 上鎖
std::cout << " Task " << i << " is done\n"; // 這裡的 i 以 1 代替
}
}
隨著 main() 函數逐漸把六個工作都塞進工作佇列 tasks、addTask() 函數也不斷喚醒執行緒,就會由三個執行緒分配六項工作。
當 main() 函數跑完迴圈(已經把六項工作加進 tasks 工作佇列,但工作可能還沒做完),進到:
// main() 函數中
return 0;
看似要結束了,但在 main() 函數中有物件實例 pool,這些要先被銷毀。
被銷毀時,呼叫解構子:
// ThreadPool 類別 public 區
~ThreadPool() {
stop();
wait();
}
首先執行 stop() 函數,給 done 上鎖時去把 done 旗標改為 true,並通知工人沒有新工作會進來,把 tasks 裡面的工作做完就好:
// ThreadPool 類別 public 區
void stop() {
{
std::lock_guard<std::mutex> lock(mtx); // 給 done 旗標上鎖
done = true; // 修改 done 旗標為 true
}
cv.notify_all(); // 告訴所有工人不會有新工作進來
}
這時旗標 done 已經被改成 true,因此不會有工人卡在這裡:
// worker() 函數中
cv.wait(lock, [this] { return !tasks.empty() || done; }); // 不會有工人在此休息
工人們依然繼續工作,此時解構子執行到 wait():
// ThreadPool 類別 public 區
void wait() {
for (auto& t : workers) {
t.join();
}
}
這裡有 join() 可確保主執行緒會等工人把工作做完,才真正銷毀 pool 物件。
等工作都做完、tasks 變成空佇列,那下面就會直接 return,離開 worker() 函數:
// worker() 函數中
if (done && tasks.empty()) {
return;
}
最後,解構子把 pool 物件都銷毀,同時運行 main() 函數中的 return 0,代表完整結束。
其中一次執行結果與說明如下:
Worker 2 starts a task
Task 1 is running <- Worker 2 搶到 Task 1
Worker 1 starts a task
Task 2 is running <- Worker 1 搶到 Task 2
Worker 0 starts a task
Task 3 is running <- Worker 0 搶到 Task 3
Task 1 is done <- Worker 2 把 Task 1 做完
Worker 2 starts a task
Task 4 is running <- Worker 2 搶到 Task 4
Task 2 is done <- Worker 1 把 Task 2 做完
Worker 1 starts a task
Task 5 is running <- Worker 1 搶到 Task 5
Task 3 is done <- Worker 0 把 Task 3 做完
Worker 0 starts a task
Task 6 is running <- Worker 0 搶到 Task 6
Task 4 is done <- Worker 2 把 Task 4 做完
Task 5 is done <- Worker 1 把 Task 5 做完
Task 6 is done <- Worker 0 把 Task 6 做完
總而言之,工作分配的結果如下:
Worker 2:Task 1 -> Task 4
Worker 1:Task 2 -> Task 5
Worker 0:Task 3 -> Task 6
每次三個工人會搶到的工作不盡相同,因此工作分配未必如上面這樣,但核心概念就是「有六個工作,交給三位工人分配」,這三個工人所在的陣列 pool 就是執行緒池。